Skip to main content

02 - LangGraph Checkpointer

前置01 篇的检查点概念。

本篇回答:LangGraph 存了什么、恢复点怎么算出来、崩溃窗口怎么处理、以及它做不到什么。

术语

  • 通道(channel):节点之间传数据的位置,可先理解为状态里的一个字段
  • 超级步(superstep):一轮并行执行。执行模型借自 Google Pregel —— 一轮内所有待执行节点并行运行,统一提交后进入下一轮
  • thread:一次会话或一个任务的标识,是检查点的主键

一、检查点的数据结构

libs/checkpoint/langgraph/checkpoint/base/__init__.py 中的定义是理解整个机制的入口:

class Checkpoint(TypedDict):
"""某一时刻的状态快照。"""

v: int
# 检查点格式版本号,当前为 1。
# 存在这个字段意味着格式会演进 —— 旧检查点必须能被新代码读懂。

id: str
# 检查点 ID。同时满足唯一且单调递增,因此可直接用于排序。
# 这样列历史时不必依赖时间戳 —— 分布式环境下时钟不可靠
# (节点间漂移、同一毫秒内产生多条)。

ts: str
# ISO 8601 时间戳。仅供展示,不参与排序。

channel_values: dict[str, Any]
# 各通道在该时刻的实际值,即"状态本身"。

channel_versions: ChannelVersions
# 各通道「当前」的版本号。每次写入该通道,版本号递增。

versions_seen: dict[str, ChannelVersions]
# 每个节点"看过"每个通道的哪个版本。
# 与 channel_versions 比对即可算出哪些节点需要执行 —— 见 1.2 节。

updated_channels: list[str] | None
# 本次检查点更新了哪些通道,用于减少下游的比对范围。

1.1 ID 单调递增带来的简化

id 同时承担主键与排序键两个职责。这个决定省掉了一整类问题:列历史记录、找父检查点、判断先后顺序全部退化成整数比较,不需要额外的序号字段,也不受时钟影响。

1.2 恢复点是算出来的,不是存下来的

这是整个机制里最值得理解的一点:LangGraph 没有"执行到第几步"这样的字段。

它存的是版本对照关系:

没有「执行到第几步」这个字段 —— 恢复点是比出来的channel_versions每个通道当前到了第几版messages :3docs :1versions_seen每个节点看过到第几版节点 A :messages 3 docs 1节点 B :messages 2比对结果谁落后了,谁就要跑节点 A 已跟上 ��—— 不执行节点 B 落后 messages —— 要执行这么设计的好处是天然支持并行和分支:不需要维护一个全局的「当前步号」,每个节点各自和通道版本比一次就知道自己该不该跑。版本号是单调递增的,直接比大小就行,不受时钟影响,也不需要额外的序号字段。
把「进度」表示成版本差而不是一个步号,是这套实现里最值得借鉴的地方 —— 它让恢复逻辑对图的拓扑变化免疫。

调度逻辑就是一句话:某节点看过的版本落后于通道当前版本,该节点就进入待执行集合。

这个设计换来三件用"程序计数器"做不到的事:

能力为什么成立
并行分支多个节点各自落后于不同通道,一轮内可全部调度
人工改状态后自动续跑手工写入某通道 → 版本号递增 → 依赖它的节点自动变为待执行
图结构变更不影响恢复依赖关系按通道计算,与节点顺序无关

可迁移的结论:为 DAG 或状态机做恢复时,记录"每个节点消费到哪个版本"比记录"全局执行到第几步"更健壮 —— 前者是局部信息、天然可并行、对拓扑变化不敏感。

1.3 全量值与增量存储

channel_values 存全量。当某个通道很大(例如一段长对话历史)时,每步存一次会造成显著的写放大 —— 即 01 篇 2.4 节描述的问题。

LangGraph 的解法是 DeltaChannel(当前标注 beta):平时只存增量写入,周期性存一次完整快照作为种子,恢复时从最近快照沿增量链向前叠加。

元数据里记录了触发快照的两个计数:

(updates, supersteps)
updates —— 有多少个超级步写过这个通道
supersteps —— 距离该通道上次快照,总共过了多少个超级步

官方注释说明了第二个计数的作用:

The supersteps bound prevents unbounded ancestor walks on threads where [the channel is rarely written]

很少被写的通道,其"最近快照"可能在几百步之前,恢复时要回溯极长的祖先链。所以除按写入次数触发快照外,还需要一个按超级步数量的强制上限兜底。

这与数据库的 WAL + checkpoint、Redis 的 AOF + RDB 是同一模式:增量存储必须配强制全量点。

二、中间写入:崩溃窗口的处理

01 篇 2.2 节指出,"步骤完成"与"检查点落盘"之间必然存在窗口。LangGraph 用 pending_writes 处理框架内部的这段窗口。

2.1 机制

class CheckpointTuple(NamedTuple):
config: RunnableConfig
checkpoint: Checkpoint
metadata: CheckpointMetadata
parent_config: RunnableConfig | None = None
pending_writes: list[PendingWrite] | None = None
# ↑ 已完成但尚未并入下一个检查点的节点输出


def put_writes(
self,
config: RunnableConfig, # 关联到哪个检查点
writes: Sequence[tuple[str, Any]], # (通道名, 值) 列表
task_id: str, # 哪个节点产生的 —— 恢复时靠它判断该节点是否已完成
task_path: str = "",
) -> None:
"""存储关联到某个检查点的中间写入。"""

方法文档里的关键词是 intermediate:一个超级步内多个节点并行执行,每个节点完成后立即单独落一次写入并带上 task_id;整轮结束后才生成新的 Checkpoint

一个超级步内部:每个节点完成就单独落一次写入节点 A 完成节点 B 完成节点 C 执行中put_writestask_id=Aput_writestask_id=B此时崩溃恢复时读到 A、B 已有写入记录只需重跑 CC 完成之后整轮结束,才生成新的 Checkpoint没有这层中间写入的话,超级步内任意一个节点崩溃,整轮三个节点都要重跑 —— 已经花掉的模型调用全部作废。代价是写入次数变多:一个超级步从写一次变成写 N+1 次。
这就是文档里 intermediate 这个词的含义:检查点的粒度不等于超级步的粒度,中间还夹了一层按节点的写入记录。

没有这个机制,整个超级步的三个节点都要重来。在"一个节点 = 一次模型调用"的场景下,这是三倍的直接成本。

2.2 它的作用范围

pending_writes 保证的是框架不会重复执行已完成的节点。它不能保证节点内部调用外部服务的幂等 —— 如果节点 A 已经发出飞书消息但 put_writes 写入失败,恢复后消息仍会重发。

外部副作用的幂等仍然由 01 篇第三节的幂等键负责,框架层解决不了。

三、接口全貌

BaseCheckpointSaver 的同步方法(每个都有 a 前缀的异步版本):

方法作用类别
get / get_tuple读取某个检查点基础
list列出历史基础
put写入检查点基础
put_writes写入中间结果基础
delete_thread删除某 thread 的全部检查点运维
delete_for_runs按 run 删除运维
copy_thread复制整个 thread能力扩展
prune清理运维
get_delta_channel_history增量通道历史(beta)能力扩展

前四个是"能存能取",后五个是"能运维、能分支"。两组的存在时间差反映了这类组件的成熟路径:先解决功能,再解决运维。

3.1 copy_thread 支撑的是调试能力

复制整个会话,是"分支"的底层依赖:任务跑到第 18 步,想验证"如果第 10 步换个方向会怎样" —— 复制一份 thread,从第 10 步改。

这是 Agent 调试的高频需求,也是 LangGraph 时间旅行功能的实现基础。

3.2 prune 与 delete_thread 不只是省存储

检查点无限增长。一个 30 步的任务若每步一个检查点、每个检查点携带完整对话历史,存储成本增长很快。

更重要的是:这两个方法同时是数据保留合规的执行点 —— 见 01 篇 2.5 节。上线前不定清理策略,三个月后会同时面对存储账单和合规问询。

四、一致性测试套件

libs/checkpoint-conformance/ 把"什么样的实现算合格的 checkpointer"写成了一套可执行测试。

# libs/checkpoint-conformance/.../validate.py
# 每一项能力对应一套独立测试
_RUNNERS = {
Capability.PUT: run_put_tests,
Capability.PUT_WRITES: run_put_writes_tests,
Capability.GET_TUPLE: run_get_tuple_tests,
Capability.LIST: run_list_tests,
Capability.DELETE_THREAD: run_delete_thread_tests,
Capability.DELETE_FOR_RUNS: run_delete_for_runs_tests,
Capability.COPY_THREAD: run_copy_thread_tests,
Capability.PRUNE: run_prune_tests,
Capability.DELTA_CHANNEL_HISTORY: run_delta_channel_history_tests,
}

async def validate(
registered: RegisteredCheckpointer,
*,
capabilities: set[str] | None = None,
...
) -> CapabilityReport:
"""对一个已注册的 checkpointer 跑验证套件。

capabilities: 指定则只测这些能力;
不指定则「自动探测」实现了哪些,只测实现了的。
"""

这是能力探测式而非全量强制式:实现了哪些就测哪些,未实现的不算失败,产出一份 CapabilityReport 能力表。

4.1 为什么这个设计值得单独学

checkpointer 是典型的"接口简单、语义微妙"组件。这些问题文档很难说清:

  • list 的排序保证是什么?
  • 并发 put 同一个 thread 会发生什么?
  • delete_thread 之后 get 应返回 None 还是报错?

语义只写在文档里,第三方实现(Redis 版、MongoDB 版、企业自研版)必然各自跑偏,出问题时无法判定是框架缺陷还是实现缺陷。

把规范写成可执行测试,是对齐语义的唯一可靠手段。 这与 Agent 安全 · MCP 授权演进中规范用 MUST / MUST NOT 收紧措辞属于同一类努力。

五、能力边界

能力支持情况
图内部状态恢复
人在环中断与恢复Interrupt 为一等公民,带 ID 可直接恢复
时间旅行 / 分支调试listcopy_thread
跨进程自动恢复
跨服务长事务❌ 仅覆盖 LangGraph 图内
定时 / 延迟执行❌ 挂起等待需自行调度
非 LangGraph 代码

5.1 最关键的缺口

检查点解决了"状态不丢",没有解决"谁把任务捡起来继续跑"。

检查点解决了「状态不丢」,没解决「谁把任务捡起来继续跑」进程崩溃检查点安静躺在 Postgres 里有人拿着 thread_id再调用一次吗没有从断点继续任务永远停在这里另外三条路线都有服务端或运行时负责发现并拉起;用 LangGraph 就得自己写一个巡检任务,扫出那些卡住的 thread_id 再调一次。
这是四条路线里唯一一处能力缺口,也是选型时最容易被忽略的一条 —— 它不体现在功能列表上,只在真的崩过一次之后才会被发现。

要补上这个缺口,需要自行实现:扫描未完成 thread、判定超时、重新拉起、并防止多实例重复拉起。最后一项最容易遗漏 —— 见 04 篇中 DBOS 用 Conductor 协调的部分。

5.2 选型判断

任务形态结论
用户主动触发、用户会回来看结果
聊天、交互式研究
LangGraph 自带的够用 —— 用户下次发消息即自然恢复
后台自动运行、无人值守
定时任务、批处理、工单流转
需要 0304 的方案

下一篇03 - Temporal 与确定性重放

← 回到 专题索引  ·  Agent Infra 板块总览